package com.wdd.flinkdemo.kafka2mysql;

import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.util.Collector;

/**
 * @author wdd
 */
public class KafkaDataHandleFunction extends ProcessFunction<String,String> {

    @Override
    public void processElement(String value, Context ctx, Collector<String> out) throws Exception {
        // TODO: 2022/4/13 做业务逻辑处理
        value = value+"modified";
        out.collect(value);
    }
}
